NOTE
2.6 Kafka Broker
Kafka Broker and Controller notes, consistency through HW, and broker-side settings related to message reliability.
This is a historical learning note and may contain outdated or incomplete understanding.
1. What Is a Broker?
- Kafka’s server side, responsible for storing messages.
- Each node in a Kafka Cluster is called a Broker. It can be regarded as an independent Kafka instance. Multiple Kafka Brokers form a Kafka Cluster.
2. Controller
2.1. What Is a Controller?
- There are one or more brokers in a Kafka cluster, and one broker is elected as the controller (Kafka Controller).
2.2. Why Have a Controller?
- Leader election: Kafka Topic. When the leader replica of a partition fails, the controller is responsible for electing a new leader replica for that partition.
- Partition reassignment: Kafka Topic. When the
kafka-topics.shscript is used to increase the number of partitions for a topic, the controller is also responsible for reassigning the partitions. - ISR change notification: when a change in a partition’s ISR set is detected, the controller is responsible for notifying all brokers to update their metadata.
2.3. Controller Election
- It depends on ZooKeeper. When each broker starts, it tries to read the broker ID from the temporary node
/controller.- If it is not -1, a controller already exists, so it gives up the election.
- Otherwise, it creates the
/controllernode. The broker that successfully creates it becomes the controller.
- ZooKeeper also has a persistent node
/controller_epoch.- It stores an integer
controller_epochvalue used to record which generation the current controller belongs to, that is, the “controller epoch”.
- It stores an integer
3. How Kafka Ensures Consistency
3.1. What It Means
- Consistency here means that the data consumers pull from any replication is the same.
- For example, suppose the Leader has offsets 1, 2, 3 and an LEO of 4, while the Follower has offsets 1, 2 and an LEO of 3.
- If a consumer can pull message 3 from the Leader and then the Leader goes down, the Follower is promoted to Leader. If message 3 cannot then be pulled from the Follower, the data is inconsistent.
3.2. How It Is Ensured
Kafka handles this situation through HW. HW = min(LEO of all replications). In the situation above, HW is 3, so consumers can only consume messages 1 and 2 and cannot consume message 3.
4. How a Broker Prevents Message Loss
4.1. Set the Number of Replicas
The replication.factor parameter sets the number of replicas for a Topic (including the leader). Each Topic can have multiple replicas, and the replicas are located on different brokers in the cluster. This means the number of replicas cannot exceed the number of brokers; otherwise, topic creation fails.
Set replication.factor >= 2, meaning at least two replicas.
4.2. Set the Number of ISRs
The min.insync.replicas parameter specifies Kafka’s minimum number of ISRs (including the leader). If the number is below this value, a Producer cannot produce messages to that Partition.
The default is 1, meaning messages can be produced as long as the leader is available. If the leader goes down, messages may be lost.
Set min.insync.replicas >= 2, meaning at least two ISRs.
4.3. Set the ACK Mechanism
The ack parameter represents the Producer’s message-send acknowledgement mechanism. There are three values:
| 0 | 1 | -1/all | |
|---|---|---|---|
| ack timing | Return ack after the Leader receives the message but before it is written to disk | Return ack after the Leader receives and writes it to disk | Return ack after the Leader receives and writes it to disk and synchronizes it to all ISRs |
| Performance | High | Medium | Low |
| Reliability | Low | Medium | High |
The default is 0. For reliability, acks = all can be set.
4.4. Set the Leader Election Mechanism
The unclean.leader.election.enable parameter specifies whether a replica outside the ISR is allowed to become leader.
The default value is true, meaning a non-ISR replica is allowed to become leader.
Set unclean.leader.election.enable = false, meaning that when the leader goes down, a leader may only be selected from the ISR.
Discussion
Sign in with GitHub to comment. Discussions are stored as GitHub Issues.View on GitHub